Skip to content

fix(grpc): forward scheduler load info for DP-aware load balancing - #1115

Closed
Kangyan-Zhou wants to merge 1 commit into
smg-project:mainfrom
Kangyan-Zhou:fix_dp_load_balance
Closed

Kangyan-Zhou wants to merge 1 commit into
smg-project:mainfrom
Kangyan-Zhou:fix_dp_load_balance

Conversation

@Kangyan-Zhou

@Kangyan-Zhou Kangyan-Zhou commented Apr 13, 2026 •

Copy link
Copy Markdown
Contributor

Summary

  • In gRPC mode, GrpcRequestManager receives BatchTokenIDOutput from the scheduler but drops the piggybacked load info (num_reqs, num_tokens per DP rank), causing total_tokens and total_requests DP load balance policies to see all-zero budgets and always pick rank 0
  • Mirror the TokenizerManager pattern: forward WatchLoadUpdateReq back to the DataParallelController when dp_size > 1

Context

In PD disaggregated serving with DP=8 decode engines, we observed severe DP rank imbalance (e.g., decoder-1 DP2 processed 93,904 tokens while DP7 processed 7). The round_robin policy distributes requests evenly by count but is blind to KV cache pressure from long-running sequences. Switching to total_tokens requires this load feedback loop to function.

Test plan

  • Deploy with --load-balance-method total_tokens on a DP=8 decode engine in gRPC mode
  • Verify DPBudget receives non-zero load updates (check scheduler debug logs)
  • Confirm DP ranks show balanced sglang_decode_sum_seq_lens under mixed-length workloads
  • Verify no regression with round_robin (load forwarding is a no-op when dp_size == 1)

🤖 Generated with Claude Code

Summary by CodeRabbit

  • Bug Fixes
    • Improved load reporting coordination in multi-instance deployments for better load distribution.

In gRPC mode, the GrpcRequestManager receives BatchTokenIDOutput from
the scheduler but does not forward the piggybacked load info
(num_reqs, num_tokens per DP rank) back to the DataParallelController.
This causes total_tokens and total_requests load balance policies to
see all-zero budgets and degenerate to always picking rank 0.

Mirror the TokenizerManager pattern: when dp_size > 1, send a
WatchLoadUpdateReq with the load info back to the scheduler input
socket so the DataParallelController can make informed dispatch
decisions.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
@coderabbitai

coderabbitai Bot commented Apr 13, 2026 •

Copy link
Copy Markdown
📝 Walkthrough

Walkthrough

Load update forwarding has been added to the batch output handler. When data parallelism is enabled (dp_size > 1), the scheduler's per-batch load information is now sent to the DataParallelController via WatchLoadUpdateReq.

Changes

Cohort / File(s) Summary
Load Update Forwarding
grpc_servicer/smg_grpc_servicer/sglang/request_manager.py
Added conditional forwarding of batch load data to DataParallelController when processing batch outputs in distributed scenarios. New import of WatchLoadUpdateReq and conditional send in _handle_batch_output method when dp_size > 1.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~12 minutes

Poem

A rabbit hops through distributed streams,
Load balancing now pursues its dreams,
Per-batch weights flow swift and bright,
Across the scheduler's watchful sight,
hippity-hop — the loads take flight! 🐰⚖️

🚥 Pre-merge checks | ✅ 3
✅ Passed checks (3 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title directly and clearly describes the main change: forwarding scheduler load information for data-parallel aware load balancing in the gRPC servicer.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands and usage tips.

@Kangyan-Zhou
Kangyan-Zhou marked this pull request as draft April 13, 2026 22:26
@mergify

mergify Bot commented Apr 13, 2026

Copy link
Copy Markdown
Contributor

Hi @Kangyan-Zhou, the DCO sign-off check has failed. All commits must include a Signed-off-by line.

To fix existing commits:

# Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-lease

To sign off future commits automatically:

  • Use git commit -s every time, or
  • VSCode: enable Git: Always Sign Off in Settings
  • PyCharm: enable Sign-off commit in the Commit tool window

@github-actions github-actions Bot added the grpc gRPC client and router changes label Apr 13, 2026

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request updates the GrpcRequestManager to forward load information to the DataParallelController when data parallelism is enabled, ensuring token-aware balancing functions correctly. The changes involve importing WatchLoadUpdateReq and adding logic to send load updates during batch output processing. A review comment suggests using the internal _send_to_scheduler helper method and awaiting the call to maintain consistency with the rest of the class and ensure proper error logging.

Comment on lines +688 to +690
if self.server_args.dp_size > 1 and batch_out.load is not None:
load_update = WatchLoadUpdateReq(loads=[batch_out.load])
self.send_to_scheduler.send_pyobj(load_update)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

To maintain consistency with the rest of the GrpcRequestManager class (see lines 368, 427, and 471), please use the _send_to_scheduler helper method instead of calling self.send_to_scheduler.send_pyobj directly. This ensures that the send operation is covered by the helper's error logging. Since _handle_batch_output is an async function, the call should be awaited to follow the established pattern in this class.

Suggested change
if self.server_args.dp_size > 1 and batch_out.load is not None:
load_update = WatchLoadUpdateReq(loads=[batch_out.load])
self.send_to_scheduler.send_pyobj(load_update)
if self.server_args.dp_size > 1 and batch_out.load is not None:
load_update = WatchLoadUpdateReq(loads=[batch_out.load])
await self._send_to_scheduler(load_update)
References
  1. If a code block is identified as duplicated across multiple functions or modules, consider refactoring to unify the logic into a shared helper function or class. This improves maintainability and reduces the chance of inconsistencies when changes are needed.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py`:
- Around line 688-690: The send of DP load updates (WatchLoadUpdateReq with
batch_out.load) from _handle_batch_output can raise and abort before put_tasks
are flushed; make this best-effort by catching exceptions around
send_to_scheduler.send_pyobj so failures won't propagate. Wrap the call that
constructs WatchLoadUpdateReq and calls self.send_to_scheduler.send_pyobj(...)
in a try/except that logs/debugs the error (non-fatal) and continues, ensuring
the subsequent flush of put_tasks still runs; keep the guard on
self.server_args.dp_size > 1 and only apply the try/except when batch_out.load
is not None.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: f2ce1e0b-ff0e-4497-9371-aab3b5f9d496

📥 Commits

Reviewing files that changed from the base of the PR and between 5319dc6 and 7858ac9.

📒 Files selected for processing (1)
  • grpc_servicer/smg_grpc_servicer/sglang/request_manager.py

Comment on lines +688 to +690
if self.server_args.dp_size > 1 and batch_out.load is not None:
load_update = WatchLoadUpdateReq(loads=[batch_out.load])
self.send_to_scheduler.send_pyobj(load_update)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🔴 Critical

Make DP load forwarding best-effort to avoid dropping batch outputs.

At Line 690, an unguarded send can raise and abort _handle_batch_output before Line 694 flushes put_tasks, which can stall/delay client-visible outputs for that batch. Handle this path as non-fatal.

Suggested fix
-        if self.server_args.dp_size > 1 and batch_out.load is not None:
-            load_update = WatchLoadUpdateReq(loads=[batch_out.load])
-            self.send_to_scheduler.send_pyobj(load_update)
+        if self.server_args.dp_size > 1 and batch_out.load is not None:
+            try:
+                await self._send_to_scheduler(WatchLoadUpdateReq(loads=[batch_out.load]))
+            except Exception as e:
+                logger.warning(f"Failed to forward DP load update: {e}")
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py` around lines 688 -
690, The send of DP load updates (WatchLoadUpdateReq with batch_out.load) from
_handle_batch_output can raise and abort before put_tasks are flushed; make this
best-effort by catching exceptions around send_to_scheduler.send_pyobj so
failures won't propagate. Wrap the call that constructs WatchLoadUpdateReq and
calls self.send_to_scheduler.send_pyobj(...) in a try/except that logs/debugs
the error (non-fatal) and continues, ensuring the subsequent flush of put_tasks
still runs; keep the guard on self.server_args.dp_size > 1 and only apply the
try/except when batch_out.load is not None.

@Kangyan-Zhou
Kangyan-Zhou marked this pull request as ready for review April 13, 2026 22:38
@mergify

mergify Bot commented Apr 13, 2026

Copy link
Copy Markdown
Contributor

Hi @Kangyan-Zhou, the DCO sign-off check has failed. All commits must include a Signed-off-by line.

To fix existing commits:

# Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-lease

To sign off future commits automatically:

  • Use git commit -s every time, or
  • VSCode: enable Git: Always Sign Off in Settings
  • PyCharm: enable Sign-off commit in the Commit tool window

@Kangyan-Zhou

Copy link
Copy Markdown
Contributor Author

Closing in favor of a new PR with correct branch naming convention and DCO signoff.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

grpc gRPC client and router changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant